{T}

穷途末路的选择:Lambda 架构

0. 引言

有些计算(二度关联图谱、复杂统计学习/机器学习模型训练)因算法复杂度过高或数据量过大,无法直接实时计算。此时需要"压箱底"的 Lambda 架构:离线与实时相结合的折中方案。本文讲解 Lambda 架构的三层设计、其改进版 Kappa 架构,并用 Flink 演示实现。

1. Lambda 架构:三层各司其职

图表渲染中…
职责存储与计算选型
批处理层处理主数据集(历史全量数据)HDFS/S3 + MapReduce/Hive/Spark;结果存 MySQL/HBase 供快速访问
快速处理层处理增量数据(新进入系统的数据)Flink/Spark Streaming/Storm;结果存 Redis 等高性能内存库
服务层合并两层结果,提供查询服务分别查询后合并返回

Lambda 架构是架构思想,各层技术选型不限定,可按实际情况组合。

主要问题:同一查询目标需为批处理层和快速处理层分别开发算法实现——相同逻辑两套代码、两个计算框架(如 Storm + Spark 并用),开发、测试、运维成本高。

2. Kappa 架构:用"流"统一编程界面

Kappa 架构的最大改进:用流计算技术取代批处理层——批处理层与快速处理层使用相同的流计算逻辑与统一的计算框架,降低开发测试运维成本。

存储也随之简化:不再将离线数据转储到 HDFS/S3 等块存储,而是按"流"的方式存于 Kafka 等流式存储系统,减少存储空间、避免数据转储、降低复杂度。

在 Flink、Spark Streaming 等流批一体计算框架与 Kafka、Pulsar 等流式存储系统的双重加持下,Kappa 已成为大数据处理的自然选择。

目标:统计"最近 3 天每种商品的销售量"。同一套代码,仅窗口类型不同:

java
DataStream counts = stream
    .map(s -> JSONObject.parseObject(s, Event.class))     // 解析事件
    .map(e -> new CountedEvent(e.product, 1, e.timestamp)) // 商品计数事件
    .assignTimestampsAndWatermarks(new EventTimestampPeriodicWatermarks())
    .keyBy("product")
    // 批处理层:滑动窗口(3 天窗口、30 分钟步长)
    .timeWindow(Time.days(3), Time.minutes(30))
    // 快速处理层:翻滚窗口(15 秒)
    // .window(TumblingEventTimeWindows.of(Time.seconds(15)))
    .reduce((e1, e2) -> merge(e1, e2));  // 窗口内聚合

批处理层与快速处理层的代码除窗口类型外完全相同——这正是 Kappa 架构的优势。

两层结果写入同一张数据库表(JdbcWriter 按层标记,ON DUPLICATE KEY UPDATE 幂等写入,并按步长取整对齐时间窗口便于合并);服务层一条 SQL 即可合并查询:

sql
SELECT product, sum(v_count) as s_count FROM (
    SELECT * FROM table_counts WHERE start=? AND end=? AND layer='batch'
    UNION
    SELECT * FROM table_counts WHERE start>=? AND end<=? AND layer='fast'
) t GROUP BY product;

4. 小结

  • Lambda 架构:不能直接实时计算时,采用"离线(批处理层)+ 实时(快速处理层)+ 合并(服务层)"的折中方案;
  • Kappa 架构:以"流"统一编程界面,同逻辑一套代码,配合流批一体框架与流式存储,显著降低复杂度;
  • Kappa 不能完全取代 Lambda:批处理算法在丰富度、成熟度、第三方工具库方面仍优于流处理;是否统一框架,还需权衡数据人员(R/Python/Spark 等熟悉工具)与开发人员各自的产出效率;
  • 最终方案应结合业务场景、技术积累、团队能力综合设计。

下一章讲解 Java 程序员的成长之路和从业方向。